Skip to content

runtime: queue-length limits for the native runtime — scope question - #35

Draft
xiaoyu-xyz wants to merge 2 commits into
ThinkFlowLab:mainfrom
xiaoyu-xyz:frontend-engine-lifecycle
Draft

xiaoyu-xyz wants to merge 2 commits into
ThinkFlowLab:mainfrom
xiaoyu-xyz:frontend-engine-lifecycle

Conversation

@xiaoyu-xyz

@xiaoyu-xyz xiaoyu-xyz commented Sep 29, 2026 •

Copy link
Copy Markdown

This PR is open as a question, not as a merge candidate. Its head (9877a6b) predates the placement revision and is out of date: it implements the queue in the frontend, which review asked to move. It is kept open because the scope question below needs an answer before the work is redone, and because its tests specify behaviour that a runtime queue would have to satisfy.

Purpose

Queue-length limits for the native runtime — whether they are wanted, and whether now is the time.

omni-runtime owns admission as SerialScheduler: FIFO, one execution unit at a time, no queue-length limit. Its README lists this explicitly:

This implements serial admission per executor within a worker process. Admitted concurrency is one. Queue-length limits, token budgets, compatibility grouping, shared processing orchestration and GPU batching remain planned.

This PR is that planned item, in the layer the architecture assigns it to. The contract it would build on landed in #34: Engine in src/runtime, with Report { depth, capacity, rejected } shaped to describe exactly this.

What #34 already has, and what would be left here

Behaviour Where it is now
the engine contract, refusal on the same reply as an answer #34, src/runtime/src/engine.rs
a worker owning a queue exercised through that contract #34, tests/runtime/engine.rs
the transport's 503 vs 504 line #34, tests/frontend/in_process.rs
refusing once capacity requests are outstanding, running or queued here
spawn returning only after load and warmup, so readiness means something here
dropping work whose caller has gone or whose budget expired, never starting it here
a failure retiring the engine instead of serving from unknown state here
draining accepted work on shutdown, and reporting when it could not finish here

The behaviour being proposed

Specified from what that implementation had to get right, and from two review findings it fixed:

  • The bound counts outstanding work, not buffer slots. The worker's channel frees its slot on dequeue; a limit that only bounded the channel let a capacity-1 worker serve two requests at once while admitting to one. Admission has to take the slot and check it under one lock, before the handover.
  • A slot is never released before it is taken. Counting after the handover lets a fast consumer decrement first; on a usize that wraps, and /health then reports a number nobody can act on.
  • A slot survives its caller. A dropped reply is how cancellation is observed; the slot is released when the work is done with, not when the caller stops listening.
  • Readiness follows warmup, not construction. spawn returns only once the model has loaded and warmed up, so a handle cannot observe a half-loaded engine.
  • A failed engine retires. A panic or an inference failure answers that caller and then stops accepting, rather than serving the next caller from state nobody can vouch for.
  • Shutdown drains and reports. Accepted work finishes within a budget; work that could not is reported rather than hidden, and admission closes at the signal rather than after the listener stops.

The question

  1. Is this in scope for omni-runtime now, or should the limits stay planned until a worker needs them? Two native workers currently run with concurrency one and no limit, and neither has reported back-pressure as a problem.
  2. If now, is the shape above the one you want — capacity as outstanding-request count with refusal at the bound, Report publishing depth/capacity/rejected, readiness tied to warmup, failure retiring the engine? Anything you would draw differently is cheaper to change before a model depends on it.
  3. Where should it live — extending SerialScheduler, or a separate admission type beside it? SerialScheduler is 40 lines and knows nothing about bodies or deadlines; a bounded queue that carries a deadline and reports depth is a different thing.

No work is pending on this: the current head is stale and would be rewritten against main and #34 if the answers say to do it. Asking first so the rewrite lands once.

@hsliuustc0106 hsliuustc0106 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Reviewed commit 1b63ed7431ea4107852600d07799a826773e3bf4.

  • [P2] Reserve queue depth before publishing the job. src/frontend/src/worker.rs:217 increments depth after try_send, but the consumer never takes the admission mutex. It can dequeue/process/decrement first, yielding depth=0 in a response and transient usize underflow in health. A standalone CPU probe against the built crate observed four zero-depth answers in 100,000 sequential immediate requests. Reserve before sending and roll back on send failure, or synchronize dequeue/accounting.
  • [P2] Start the drain budget when shutdown is signaled. src/frontend/src/main.rs:68 awaits graceful HTTP shutdown before stop_accepting/run_until_drained. In-flight upload/response handlers can consume their entire request budget before the drain timer even begins, and uploads can still submit during this interval. A real binary with DRAIN_MS=50 and TIMEOUT_MS=3000 took 2.917 seconds to exit after SIGTERM with one partial upload. Stop admission on the signal and bound the complete shutdown sequence with the drain deadline.
    All 31 ordinary tests plus 1 doctest passed; these probes expose gaps in the existing suite.

@xiaoyu-xyz
xiaoyu-xyz force-pushed the frontend-engine-lifecycle branch from 1b63ed7 to 3219451 Compare September 30, 2026 00:55
@xiaoyu-xyz xiaoyu-xyz changed the title feat(frontend): own the engine in a worker thread with a bounded queue [stacked on #34] feat(frontend): own the engine in a worker thread with a bounded queue Sep 30, 2026
@xiaoyu-xyz

xiaoyu-xyz commented Sep 30, 2026 •

Copy link
Copy Markdown
Author

Superseded — see the current status comment.

This comment is kept only so the thread does not lose a reply that was read. Two of its claims did not survive re-checking, and I am not leaving them standing as if they had:

  1. It said the depth test "fires within the first hundred submissions" against the released-before-reserved ordering. It does not. That observation came from an earlier version of the test, before I changed the fake; the test as it exists passes against that ordering too. Its own doc comment says so, and the current comment states what the test actually covers.
  2. It reported "34 passed" and a set of timings measured on a head that has since been rebased twice. Those numbers were current when written and are not now — the same head now has far more tests, because main keeps adding them.

Nothing else in it was withdrawn — the two P2 findings were fixed as described, and the shutdown and capacity findings that came afterwards are covered in the current comment.

@Levius-Fubuki Levius-Fubuki left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Reviewed 3219451e57ab7ad96fc92ae47d4636fe404f9c7b, including the incremental changes on #34 and the two previous findings.

The depth reservation now precedes publication and rolls back on both failed-send paths, addressing the original counter-underflow mechanism. However, the shutdown follow-up introduces a reproducible service-lifetime regression. Changes are requested for the three issues below.

Formatting passed. Clippy, workspace tests and release build passed on Linux/Rust 1.98.1 after rebuilding in a dedicated target directory (an initial shared-target run had stale cross-checkout artifacts and is not treated as a PR failure). Independent probes reproduce the no-signal exit and invalid-mode fallback on Linux and macOS; a gated CPU worker reproduces the capacity breach. No GPU is needed for these cases. This remains a draft and also needs its existing merge conflicts resolved.

Comment thread src/frontend/src/main.rs Outdated
// Responses in flight get the drain budget to finish. A handler still stuck on a slow
// upload at the deadline is not going to finish, and holding the process open for it
// would mean the budget bounded nothing.
let deadline = Instant::now() + running.drain_budget();

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P1] Start the drain deadline on the shutdown signal, not at startup

This deadline is created as soon as the listener starts and wraps the entire engine::serve future. Consequently every native process exits after the drain duration even if no SIGTERM/SIGINT is sent (10 seconds by default). With OMNI_SYSTEMONE_ENGINE=passthrough, OMNI_SYSTEMONE_DRAIN_MS=200 and OMNI_JEV_BIND=127.0.0.1:0, the release binary exited 0 after 0.206 seconds on Linux, logging shutdown exceeded its drain budget; the probe sent no signal. Wait for the signal while serving normally, then close admission and compute one absolute deadline for graceful HTTP shutdown plus worker drain. Add a test that remains healthy beyond the drain duration before sending any signal.

Comment thread src/frontend/src/worker.rs Outdated
// answered and released — counting afterwards would release a slot that had not
// been taken yet, which reads as a wrong depth and wraps the counter. The window is
// small but real, and `depth_is_reserved_before_the_job_is_published` caught it.
self.depth.fetch_add(1, Ordering::AcqRel);

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] Enforce the documented limit on all outstanding requests

This increment reserves accounting but never checks the configured outstanding-request limit; only mpsc::try_send enforces capacity. Once the worker dequeues a job, the channel slot is free even while inference still runs. A deterministic probe with capacity 1 holds the first inference behind a gate, waits until it has started, then submits a second job: it is accepted and report().depth becomes 2. Both Options::queue_capacity and the README promise refusal once that many accepted requests are unfinished. Check/reserve the outstanding count under the admission lock (with rollback on send failure), and cover a running-plus-queued case instead of relying on whether the consumer has happened to dequeue.

Comment thread src/frontend/src/main.rs Outdated

async fn run() -> Result<(), BoxError> {
let bind = setting("OMNI_JEV_BIND", Config::DEFAULT_BIND.to_owned())?;
if env::var("OMNI_SYSTEMONE_ENGINE").as_deref() != Ok("passthrough") {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] Reject unsupported engine names instead of silently selecting forwarding

Every value other than passthrough enters forwarding mode here. With OMNI_SYSTEMONE_ENGINE=typo-engine, the built binary starts and logs forwarding to http://127.0.0.1:8000/, although the README says unknown modes are startup errors. A configuration typo can therefore route inference to a different backend instead of failing visibly. Match supported values explicitly, preserve the documented unset/default behavior, and return a configuration error for unknown values.

@Levius-Fubuki

Copy link
Copy Markdown
Collaborator

Rechecked a9f7b83: the no-signal exit, capacity limit, and unknown-engine handling are fixed. Shutdown still exceeds the drain budget with an open, incomplete upload: with a 200 ms drain budget and 3 s request timeout, SIGTERM took ~2.85 s to exit in two runs. The new test drops the connection before signaling, so it misses this case. Please keep the connection open in the test and bound the entire HTTP/worker shutdown sequence.

@xiaoyu-xyz
xiaoyu-xyz force-pushed the frontend-engine-lifecycle branch 2 times, most recently from a0e86f4 to c88dc14 Compare October 4, 2026 04:59
RFC ThinkFlowLab#14 puts the request lifecycle in src/frontend/ and the model behind "a small
engine interface", and #3's acceptance criteria depend on that interface existing.
Neither half is on main today.

This adds the interface and the HTTP service behind it: the Engine trait, its
readiness/answer/error types, and engine::app serving /v1/systemone and /health
with a request budget measured from the headers, a body limit, and one status
code per failure mode.

Nothing parses the decision envelope, so a field this crate has never heard of
survives and a model's error text cannot produce invalid JSON. submit takes owned
bytes and returns a channel, which is the only shape that works for a worker
thread with a queue behind it; the reply carries the queue depth to publish, so
x-queue-depth is measured by the engine rather than guessed by the transport.

No model, no GPU, and no change to main.rs or the forwarding path: the thread,
the bounded queue and the drain that plug into this boundary come next.
@xiaoyu-xyz
xiaoyu-xyz force-pushed the frontend-engine-lifecycle branch from c88dc14 to 63b641c Compare October 4, 2026 14:37
@xiaoyu-xyz

Copy link
Copy Markdown
Author

@Levius-Fubuki — the shutdown case you flagged is fixed, and you were right about why my test missed it.

Your review was against a9f7b83; the fix landed after it. I read your comment with the review, but the review itself was posted against that older commit, so it looked like the bound had already been addressed when it had not. That is on me for not checking which head each comment was anchored to.

The measurement you reported reproduces here on a9f7b83: 2.60 s, 2.61 s with DRAIN_MS=200 and TIMEOUT_MS=3000 — the request budget, not the drain budget, was bounding the shutdown.

Why the wait cannot be signalled away. I instrumented the two halves separately: the shutdown future finished at 1.01 s, and serve returned at 3.41 s. The difference is that connection's request budget. axum's graceful shutdown waits for open connections for as long as they are open, and dropping the future at the deadline does not close them. So the budget has to be enforced by no longer waiting on the transport, not by asking it to stop.

My test was the reason I did not see it. It hung up before signalling, and the bug only appears while the connection is open — I wrote a test for a bug and then arranged the precondition away. It now holds stalled open across the signal and drops it only after child.wait(), and it fails against a9f7b83:

Measured with the connection held open through the whole signal-to-exit window:

Budget Request timeout a9f7b83 63b641c
DRAIN_MS=50 3000 2.61 s, 2.60 s 339 ms, 345 ms, 341 ms
DRAIN_MS=200 3000 2.62 s, 2.63 s 506 ms, 504 ms, 500 ms

The bound now runs from the signal and only from the signal: the shutdown future closes admission, starts the budget, drains the worker, and serves a bounded grace period for responses already in flight; the outer wait ends on whichever comes first — serve resolving, the grace window ending, or the budget. With no signal sent, the process keeps serving, which the_service_stays_up_beyond_the_drain_budget_without_a_signal covers.

Two wrong attempts before this one, both of which looked right: timing the serve future from startup re-introduced the no-signal exit, and a Notify to wake the outer wait stores a permit, so it fired immediately and skipped the grace period. An oneshot does not.

Also rebased onto current main (6d92519) while I was here, so mergeable is true; the README conflict was structural and is described on #34. fmt and clippy -D warnings clean, release build passes, and the shutdown tests pass on 63b641c.

@Levius-Fubuki

Copy link
Copy Markdown
Collaborator

Rechecked 63b641c: open-upload shutdown is improved, but still exceeds the drain budget. From SIGTERM, 50 ms / 200 ms budgets took ~304 ms / ~454 ms, each reproduced three times. The outer timeout(DRAIN_GRACE, ...) adds another 250 ms wait. Please use the same absolute deadline throughout and tighten the regression assertion to catch this overrun.

@xiaoyu-xyz
xiaoyu-xyz force-pushed the frontend-engine-lifecycle branch from 63b641c to 9876a7b Compare October 4, 2026 15:48
@xiaoyu-xyz

Copy link
Copy Markdown
Author

@Levius-Fubuki — one absolute deadline now, and the regression test is tight enough to stay red on the additive version.

What changed. The deadline is taken once, at the signal, and every step after it is bounded by that instant rather than by a fresh period:

  • The worker drains until that instant (run_until_drained(deadline)), as before.
  • The transport's close window is now deadline.min(now + CLOSE_GRACE) — clamped to the deadline instead of deadline + grace. The earlier version slept its grace period after the drain and then handed the select a fresh 250 ms, which is two additive periods for one budget.
  • The outer wait polls the server future until that instant rather than awaiting it once, so it ends when the transport is done instead of always lasting the window.

That also fixed a regression I introduced while bounding this: because the window was the whole budget, a binary with DRAIN_MS=10000 took its full 10 s to exit, and binary_serves_a_linked_engine_and_drains_on_sigterm — which has been in the branch since before this review — became flaky and then slow. It is back to ~1.4 s for the whole file.

A regression test for the overrun itself. shutdown_does_not_wait_out_a_budget_it_does_not_need uses a 5 s budget with an idle keep-alive connection held open and asserts exit under 1.5 s. An additive grace would spend budget + grace; a window without a bound spends the whole budget. The existing shutdown_is_bounded_by_the_drain_budget_not_the_request_budget keeps the connection open across the signal (you were right that dropping it first is what hid this) and now asserts under 450 ms rather than a loose second.

One thing I have to be straight about. I cannot reproduce the numbers you measured on 63b641c. With a probe that holds an incomplete upload open, and separately an idle keep-alive connection, 63b641c and the current head both exit in ~240 ms at a 200 ms budget and ~290 ms at 1 s and above — identical, three runs each. I also tried the exact configuration from your earlier comment (DRAIN_MS=50, TIMEOUT_MS=3000) and got 68–76 ms on both. So either your probe differs from mine in a way I have not found, or something in your environment adds the 60–160 ms.

What I can state as measured, with the connection held open through the shutdown:

drain budget this head 63b641c
50 ms 67–89 ms 68–76 ms
200 ms 233–260 ms 229–241 ms
1000 ms 253–293 ms 286–293 ms
10 000 ms 255–295 ms 285–297 ms

And the structural claim, which is what your instruction asked for and what the source now shows: there is one deadline, taken at the signal; the worker drains to it; the transport's window is clamped to it; nothing adds a period measured afterwards. If your probe still shows an overrun on 9876a7b, its output would tell me what mine is missing — I would rather fix that than argue from my own harness.

Verified on a clean checkout of 9876a7b, changed files hashed against the pushed branch: fmt and clippy -D warnings clean, 57 tests passing across six consecutive runs from a fresh build, release build passes.

@xiaoyu-xyz

xiaoyu-xyz commented Oct 5, 2026 •

Copy link
Copy Markdown
Author

Superseded by #34 (reopened with the placement revision; #91 was a duplicate and is closed). Closing this branch.

The placement finding on #34 applies here more directly than anywhere: this is the worker-thread queue and lifecycle that the finding is about. The contract it would implement now lives in src/runtime, and omni-runtime owns admission already as SerialScheduler, with queue-length limits listed in its README as planned.

A correction to what I first wrote here. I said the tests worth keeping from this branch are in #34. That was too broad, and I checked rather than leave it standing:

This branch had In #34 now
queue capacity, depth accounting, refusal does not count yes — tests/runtime/engine.rs
the line between 503 and 504 yes — tests/frontend/in_process.rs
spawn returns only after load and warmup no
a startup failure is reported and leaves nothing ready no
cancelled and expired requests never reach the engine no
a panic is answered, then retires the engine no
drain reports when it could not finish no
a refused request cannot leave depth permanently high no

So nine of those thirteen are not in #34; they belong to the queue and lifecycle, and they went out with it. Two of the seven rows above were also findings from review that were fixed here first — the capacity definition and the depth reservation — and the second has a known limit that has not changed: the underflow interleaving was never reproduced by my own test, only by the reviewer's probe.

What this branch leaves behind, then, is a specification of the behaviour a runtime queue would need: readiness that follows warmup, a bound that never counts a slot it did not take, work that is dropped rather than started when its caller has gone or its budget expired, a failure that retires the engine instead of serving from unknown state, and a drain that reports when it could not finish. If queue-length limits are picked up, those are the cases to cover, and the branch is here to read.

@xiaoyu-xyz xiaoyu-xyz closed this Oct 5, 2026
@xiaoyu-xyz xiaoyu-xyz reopened this Oct 5, 2026
@xiaoyu-xyz xiaoyu-xyz changed the title feat(frontend): own the engine in a worker thread with a bounded queue runtime: queue-length limits for the native runtime — scope question Oct 5, 2026
@xiaoyu-xyz

Copy link
Copy Markdown
Author

The question this PR is open for, in one place, so it can be answered without reading the stale diff.

omni-runtime's README lists queue-length limits as planned. This branch implemented them in the frontend, which review correctly rejected — they belong in the runtime, and #34 now carries the contract they would build on. Before rewriting the implementation against #34, three things need a decision:

  1. In scope now, or still planned? Two native workers run with concurrency one and no limit today, and neither has reported back-pressure as a problem. If a limit is not needed yet, this stays planned and I will not spend the rewrite.

  2. If now, is this the shape? Refuse once capacity requests are outstanding — running or queued, not merely buffered, since the channel frees its slot on dequeue and a buffer-only limit let a capacity-1 worker admit two. A slot is taken before the handover and rolled back if the handover fails, so nothing is deducted that was never counted. A dropped reply does not release a slot early; it is released when the work is done with. Readiness follows warmup rather than construction. A failed or panicking engine answers its caller and then retires instead of serving from unknown state. Shutdown closes admission at the signal and reports a drain it could not finish.

  3. Where does it live? SerialScheduler is 40 lines and knows nothing about bodies or deadlines; a bounded queue carrying a deadline and reporting depth is a different thing. Extend it, or add a type beside it?

Report { depth, capacity, rejected } in #34 was shaped for this, so the reporting half is already agreed unless you would draw it differently.

Nothing is pending on an answer — this branch's head is stale and would be rewritten against main and #34 if the answer is to proceed. Asking first so that rewrite happens once.

@Levius-Fubuki

Copy link
Copy Markdown
Collaborator

Rechecked 9876a7b: the extra shutdown wait is fixed; 50/200 ms budgets exited in ~51/201 ms (three runs each). One workspace run failed the cancellation test's depth assertion; 20 isolated suite reruns and a full workspace rerun passed (57 passed, 9 ignored). The reply can arrive before slot release, so please synchronize that assertion. This remains a stale frontend implementation pending the runtime-scope decision.

@xiaoyu-xyz
xiaoyu-xyz force-pushed the frontend-engine-lifecycle branch from 9876a7b to cbbeedf Compare October 5, 2026 12:58
WARNING: this head is stale and is not a merge candidate. It implements the
queue in the frontend, which review asked to move to the worker-side runtime,
and it predates the contract it would build on (`Engine` now lives in
src/runtime). It is kept because the queue-length question in the PR body needs
an answer before the work is redone against current `main`, and because the
tests here specify behaviour a runtime queue would have to satisfy.

What it is, for reading:

- worker.rs: one thread owns the engine, a bounded queue owns admission, and
  readiness, cancellation, expiry and drain are the transport's. A slot is taken
  before a job is handed over and rolled back if the handover fails; a job whose
  caller has gone or whose budget expired is dropped instead of started; a panic
  in model code answers its caller and retires the engine.
- passthrough.rs: a stand-in engine that echoes requests, so the path can run
  with no model and no GPU.
- worker_lifecycle.rs: the cases above, each with a failure mode behind it.

The last change synchronizes three depth assertions against the worker's own
reply ordering. The worker sends a reply before it releases the slot that reply
held, so reading depth immediately after `reply` could still see the slot in
use and fail about once in a full workspace run. Depth only falls as the worker
finishes, so the checks now wait for zero with a deadline. I could not reproduce
the window locally in twelve instrumented runs, so the reason for the change is
the ordering in the code and the reported failure, not my own observation.
@xiaoyu-xyz
xiaoyu-xyz force-pushed the frontend-engine-lifecycle branch from cbbeedf to c423100 Compare October 5, 2026 12:59
@xiaoyu-xyz

Copy link
Copy Markdown
Author

Synchronized on c423100. The assertion was reading depth immediately after a reply, and the worker sends that reply before it releases the slot it was holding, so the read could land in between.

The three depth checks now wait for zero rather than sampling once. Depth only falls as the worker finishes, so waiting is what checks the property actually being claimed — every accepted request releases its slot — and a deadline still fails the test if some request never does. I converted the same shape in the capacity test and the storm test while I was there.

One thing I have to be straight about: I could not reproduce the window locally. I instrumented the read to print the depth at the moment the reply arrived and ran it twelve times; it read zero every time. So the reason for the change is the ordering in the source — job.reply.send(...) before the slot's guard drops — plus your reported failure, not my own observation of it firing. If your run failed once in a workspace pass, that is a more sensitive probe than mine and I am treating it as the measurement.

After the change: 40 isolated runs of that test, 15 runs of the whole suite, and 12 full workspace runs, all passing (57 passed, 9 ignored).

On the stale-implementation note: agreed, and the branch says so now — the commit message opens by stating it is not a merge candidate and why, and the PR body asks the scope question rather than proposing the diff. No work is pending on it.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants